ZeroMQ 代理心跳检测:让 proxy 感知 server 存活
用 zmq_proxy() 搭建异步客户端/服务器模式的代理时,有一个容易被忽视的坑:当 server 全部退出或网络中断后,管道里积压的未消费消息会让 proxy 产生高额 CPU 开销。解法是让 proxy 能感知 server 是否存活——server 数量归零时及时关闭代理。本文介绍用 ZMTP 心跳与 socket monitor 事件监听组合实现这一检测,附完整可运行的 C++ 示例。适合正在用 libzmq 搭建 ROUTER/DEALER 代理的开发者。
异步客户端/服务器模式与 zmq_proxy
ZeroMQ 官方指南的异步客户端/服务器模式用代理连接两侧:多个 client 与多个 server 不直接通信,中间由一个 proxy 转发。
代理本身由 zmq_proxy() 实现:
int zmq_proxy(const void *frontend, const void *backend, const void *capture); 几个关键行为(依据 libzmq 官方文档):
zmq_proxy()在当前线程中运行,仅在当前上下文关闭时返回(此时返回 -1,errno 为 ETERM 等);- 代理将前端连接到后端,数据从前端流向后端只是概念上的——代理是完全对称的,前后端之间没有技术差异;
- 如果 capture 参数不为 NULL,代理会把前后端收到的所有消息再发送到 capture 套接字,capture 应为
ZMQ_PUB、ZMQ_DEALER、ZMQ_PUSH或ZMQ_PAIR类型,常用于旁路观察流量; - 调用前需要先设置必要的套接字选项,并完成 bind 或 connect。
为什么需要检测 server 存活
zmq_proxy() 在管道中存在积压的未消费消息时,会持续空转并产生高额 CPU 开销。导致积压的典型场景有两个:
- server 进程退出(正常退出或崩溃);
- 网络中断,server 无法触达但连接尚未被判定失效。
这两种场景需要不同的检测手段,区别在于 TCP 连接的回收机制:
- 进程退出时:TCP 连接信息由内核维护。进程退出后,内核回收其所有 TCP 连接资源,主动发送第一次挥手的 FIN 报文,后续挥手过程也在内核完成,不需要进程参与。也就是说,即使 server 进程已经不存在,两端仍能完成 TCP 四次挥手。此时只需用 socket monitor 监听
ZMQ_EVENT_DISCONNECTED事件,即可感知断开。 - 网络中断时:TCP 连接仍然存在,对端崩溃、拔线、断网都不会触发
ZMQ_EVENT_DISCONNECTED。这时需要在应用层发送心跳包,靠超时判定 server 是否存活。
ZMTP 心跳:三个套接字选项
zmq_setsockopt() 中与心跳相关的选项有三个:
| 选项 | 语义 | 依赖与生效条件 |
|---|---|---|
ZMQ_HEARTBEAT_IVL | 心跳发送间隔(毫秒)。大于 0 时,每隔 IVL 毫秒发送一次 PING ZMTP 命令 | 独立生效,默认 0(不发心跳) |
ZMQ_HEARTBEAT_TIMEOUT | 发送 PING 后未收到任何流量时,判定连接超时的等待时长(毫秒) | 仅在 IVL 已设置且大于 0 时有效;默认取 IVL 值 |
ZMQ_HEARTBEAT_TTL | 告知远端的心跳超时(毫秒)。远端在 TTL 内未收到任何流量则断开连接 | 依赖 IVL;内部向下取整到 0.1 秒,小于 100 的值无效,最大 6553599 |
两个容易误读的细节(均为官方文档明确说明):
- 取消 TIMEOUT 超时不要求收到 PONG——任何收到的流量都会取消本次超时计时;
- TTL 是设置在远端生效的:本端通过该选项告知对端「如果 TTL 时间内没有收到任何流量,就由你那边超时断开」。
在代理的后端套接字上添加心跳(需要在 bind 之前设置,套接字选项只影响之后建立的连接):
_backend.setsockopt(ZMQ_HEARTBEAT_IVL, 1000);
_backend.setsockopt(ZMQ_HEARTBEAT_TIMEOUT, 5000);
_frontend.bind("tcp://*:5570");
_backend.bind("tcp://*:5571");
_capture.bind("inproc://capture"); 这样 proxy 会对连接到 backend 的每一个 TCP 连接发送心跳。设置 socket monitor 之后,心跳超时导致的断开同样会触发 ZMQ_EVENT_DISCONNECTED——实测在局域网内两台机器分别启动 proxy 与 server,断开网络连接后,达到超时时间即触发该事件。用 tcpdump -i lo -vnnX port 5571 抓包,可以看到 proxy 与 server 之间一秒一次的 PING/PONG 心跳。
用 socket monitor 捕获断开事件
int zmq_socket_monitor(void *socket, char *endpoint, int events); zmq_socket_monitor() 允许应用跟踪某个套接字上的连接事件。每次调用会创建一个 ZMQ_PAIR 套接字并绑定到指定的 inproc:// 端点(监控端点只支持 inproc 传输;被监控的套接字本身须使用面向连接的传输,如 tcp/ipc)。应用自己创建一个 ZMQ_PAIR 套接字连接到该端点,即可收到事件。events 参数是事件位掩码,ZMQ_EVENT_ALL 表示全部事件。
每个事件以两帧发送:
- 第一帧:事件编号(16 位)+ 事件值(32 位,按事件类型提供附加数据,例如
ZMQ_EVENT_DISCONNECTED的值是底层网络套接字的文件描述符); - 第二帧:受影响端点的字符串。
解析两帧的官方示例代码:
static int get_monitor_event(void *monitor, int *value, char **address)
{
// 第一帧:事件编号(16位)+ 事件值(32位)
zmq_msg_t msg;
zmq_msg_init(&msg);
if (zmq_msg_recv(&msg, monitor, 0) == -1) return -1;
assert(zmq_msg_more(&msg));
uint8_t *data = (uint8_t *)zmq_msg_data(&msg);
uint16_t event = *(uint16_t *)(data); // 事件编号
if (value) *value = *(uint32_t *)(data + 2); // 事件值
// 第二帧:受影响端点的字符串
zmq_msg_init(&msg);
if (zmq_msg_recv(&msg, monitor, 0) == -1) return -1;
assert(!zmq_msg_more(&msg));
if (address) {
uint8_t *data = (uint8_t *)zmq_msg_data(&msg);
size_t size = zmq_msg_size(&msg);
*address = (char *)malloc(size + 1);
memcpy(*address, data, size);
(*address)[size] = 0;
}
return event;
} 事件循环中维护一个 server 计数:握手成功(ZMQ_EVENT_HANDSHAKE_SUCCEEDED)时加一,断开(ZMQ_EVENT_DISCONNECTED)时减一,归零时触发关闭代理:
// 套接字监控仅适用于 inproc://
int rc = zmq_socket_monitor(_backend, "inproc://monitor", ZMQ_EVENT_ALL);
assert(rc == 0);
zmq::socket_t server_mon(_ctx, ZMQ_PAIR);
std::thread th1([&]() {
server_mon.connect("inproc://monitor");
int event = 0;
while (event != ZMQ_EVENT_MONITOR_STOPPED)
{
event = get_monitor_event(server_mon, NULL, NULL);
switch (event)
{
case ZMQ_EVENT_HANDSHAKE_SUCCEEDED:
++_serverCnt;
break;
case ZMQ_EVENT_DISCONNECTED:
--_serverCnt;
if (_serverCnt == 0)
{
// 关闭 proxy,注意资源回收
}
break;
default:
break;
}
}
}); 一个实现上的提醒:示例代码里关闭 proxy 用的是 pthread_cancel(),这有可能导致 ZeroMQ 上下文崩溃,仅作演示。生产实现建议改用 zmq_proxy_steerable() 代替 zmq_proxy(),向其 control 端点发送 TERMINATE 指令来优雅终止代理。
完整示例
以下代码实现多对多的异步客户端/服务器模式:proxy 统计握手成功的 server 数量,归零后关闭代理。三段代码可分别编译运行。
proxy
#include <czmq.h>
#include <atomic>
#include <condition_variable>
#include <iostream>
#include <thread>
#include <zmq.hpp>
class Proxy
{
public:
Proxy()
: _ctx(1),
_frontend(_ctx, ZMQ_ROUTER),
_backend(_ctx, ZMQ_DEALER),
_capture(_ctx, ZMQ_DEALER),
_eventCollect(_ctx, ZMQ_PAIR),
_msgCollect(_ctx, ZMQ_DEALER)
{
}
void init()
{
try
{
// server 归零后异步关闭 proxy
std::thread t([&]() {
while (true)
{
std::unique_lock<std::mutex> lck(_mtx);
_cond.wait(lck);
// 等待 server 重启窗口
std::this_thread::sleep_for(std::chrono::seconds(3));
if (_serverCnt == 0)
{
// 生产环境建议用 proxy_steurable() + TERMINATE 指令
pthread_cancel(_proxyThread.native_handle());
break;
}
}
});
t.detach();
// 套接字监控仅适用于 inproc://
int rc = zmq_socket_monitor(_backend, "inproc://monitor", ZMQ_EVENT_ALL);
assert(rc == 0);
std::thread th1([&]() {
_eventCollect.connect("inproc://monitor");
int event = 0;
while (event != ZMQ_EVENT_MONITOR_STOPPED)
{
event = get_monitor_event(_eventCollect, NULL, NULL);
switch (event)
{
case ZMQ_EVENT_HANDSHAKE_SUCCEEDED:
++_serverCnt;
break;
case ZMQ_EVENT_DISCONNECTED:
--_serverCnt;
if (_serverCnt == 0)
{
_cond.notify_one();
}
break;
default:
break;
}
}
});
_monitorThread.swap(th1);
_frontend.setsockopt(ZMQ_SNDHWM, 10000);
_backend.setsockopt(ZMQ_HEARTBEAT_IVL, 1000);
_backend.setsockopt(ZMQ_HEARTBEAT_TIMEOUT, 5000);
_frontend.bind("tcp://*:5570");
_backend.bind("tcp://*:5571");
_capture.bind("inproc://capture");
// 捕获消息的线程:旁路统计经过 proxy 的流量
std::thread th2([&]() {
_msgCollect.connect("inproc://capture");
int cnt = 0;
while (true)
{
zmq::message_t identity;
zmq::message_t msg;
_msgCollect.recv(&identity);
_msgCollect.recv(&msg);
std::cout << "count:" << ++cnt << std::endl;
}
});
_captureThread.swap(th2);
} catch (std::exception &e)
{
std::cerr << e.what();
}
}
void join()
{
_proxyThread.join();
_frontend.close();
_backend.close();
_capture.close();
pthread_cancel(_captureThread.native_handle());
_captureThread.join();
_msgCollect.close();
_monitorThread.join();
_eventCollect.close();
_ctx.close();
}
void run()
{
try
{
std::thread th([&]() {
zmq::proxy(static_cast<void *>(_frontend), static_cast<void *>(_backend),
static_cast<void *>(_capture));
});
_proxyThread.swap(th);
std::cout << "proxy running..." << std::endl;
} catch (std::exception &e)
{
std::cerr << e.what();
}
}
private:
int get_monitor_event(void *monitor, int *value, char **address)
{
zmq_msg_t msg;
zmq_msg_init(&msg);
if (zmq_msg_recv(&msg, monitor, 0) == -1) return -1;
uint8_t *data = (uint8_t *)zmq_msg_data(&msg);
uint16_t event = *(uint16_t *)(data);
if (value) *value = *(uint32_t *)(data + 2);
zmq_msg_init(&msg);
if (zmq_msg_recv(&msg, monitor, 0) == -1) return -1;
if (address)
{
uint8_t *data = (uint8_t *)zmq_msg_data(&msg);
size_t size = zmq_msg_size(&msg);
*address = (char *)malloc(size + 1);
memcpy(*address, data, size);
(*address)[size] = 0;
}
return event;
}
private:
zmq::context_t _ctx;
zmq::socket_t _frontend;
zmq::socket_t _backend;
zmq::socket_t _capture;
zmq::socket_t _eventCollect;
zmq::socket_t _msgCollect;
std::atomic<int> _serverCnt{0};
std::thread _proxyThread;
std::thread _captureThread;
std::thread _monitorThread;
mutable std::mutex _mtx;
std::condition_variable _cond;
};
int main()
{
Proxy proxy;
proxy.init();
proxy.run();
proxy.join();
} server
#include <czmq.h>
#include <thread>
#include <zmq.hpp>
class Server
{
public:
Server() : _ctx(1), _serverSocket(_ctx, ZMQ_DEALER) {}
void start()
{
_serverSocket.connect("tcp://localhost:5571");
try
{
while (true)
{
zmq::message_t identity;
zmq::message_t msg;
zmq::message_t copied_id;
_serverSocket.recv(&identity);
_serverSocket.recv(&msg);
// 每个请求回复 5 条,验证多对多转发
for (int reply = 0; reply < 5; ++reply)
{
copied_id.copy(&identity);
_serverSocket.send(copied_id, ZMQ_SNDMORE);
char response_string[50] = {};
sprintf(response_string, "%s, response #%d", msg.to_string().c_str(), reply);
_serverSocket.send(response_string, strlen(response_string), ZMQ_DONTWAIT);
std::this_thread::sleep_for(std::chrono::milliseconds(1));
}
}
} catch (std::exception &e)
{
}
}
private:
zmq::context_t _ctx;
zmq::socket_t _serverSocket;
};
int main()
{
Server ser;
ser.start();
} client
#include <czmq.h>
#include <iostream>
#include <thread>
#include <zmq.hpp>
#define within(num) (int)((float)((num)*random()) / (RAND_MAX + 1.0))
class Client
{
public:
Client() : _ctx(1), _clientSocket(_ctx, ZMQ_DEALER) {}
void start()
{
char identity[10] = {};
sprintf(identity, "%04X-%04X", within(0x10000), within(0x10000));
_clientSocket.setsockopt(ZMQ_IDENTITY, identity, strlen(identity));
_clientSocket.setsockopt(ZMQ_RCVHWM, 5000);
_clientSocket.connect("tcp://localhost:5570");
zmq::pollitem_t items[] = {{_clientSocket, 0, ZMQ_POLLIN, 0}};
int request_nbr = 0;
try
{
std::thread th([&]() {
while (true)
{
zmq::poll(items, 1, 1);
if (items[0].revents & ZMQ_POLLIN)
{
s_dump(_clientSocket);
}
}
});
int cnt = 100;
do
{
char request_string[16] = {};
sprintf(request_string, "request #%d", ++request_nbr);
std::cout << "send " << request_string << std::endl;
_clientSocket.send(request_string, strlen(request_string), ZMQ_DONTWAIT);
std::this_thread::sleep_for(std::chrono::seconds(1));
} while (--cnt);
th.join();
} catch (std::exception &e)
{
}
}
private:
zmq::context_t _ctx;
zmq::socket_t _clientSocket;
};
int main()
{
Client cli;
cli.start();
} 要点
- server 进程退出由内核完成 TCP 四次挥手,socket monitor 的
ZMQ_EVENT_DISCONNECTED就能感知;网络中断必须靠 ZMTP 心跳超时判定,两种场景可以统一收敛到 DISCONNECTED 事件上处理。 - 心跳三选项中,TIMEOUT 依赖 IVL;取消超时不要求收到 PONG,任何流量都算数;TTL 是告知远端的超时,向下取整到 0.1 秒。
- 心跳选项须在 bind/connect 之前设置,只影响之后建立的连接。
- 关闭代理优先使用
zmq_proxy_steerable()+TERMINATE指令,避免对 ZeroMQ 线程使用pthread_cancel()。